Skip to content

[SPARK-59624][PYTHON] Move element-wise UDF conversion behind Arrow batch interfaces - #58903

Closed
Spenserrrr wants to merge 5 commits into
apache:masterfrom
Spenserrrr:arrow-worker-batch-interface
Closed

Spenserrrr wants to merge 5 commits into
apache:masterfrom
Spenserrrr:arrow-worker-batch-interface

Conversation

@Spenserrrr

@Spenserrrr Spenserrrr commented Sep 18, 2026 •

Copy link
Copy Markdown
Contributor

What changes were proposed in this pull request?

  • Add a general ArrowBatchTransformer.concat_batches classmethod with compatibility for PyArrow versions before 19.0.0.
  • Keep element-wise list flattening and shape reconstruction local to the worker while wrapping flattened inputs in RecordBatches for batch-level pandas/Arrow conversion.
  • Reuse the shared concat helper in the grouped aggregate and window Arrow handlers migrated by #58956, while preserving existing UDF contracts.

Why are the changes needed?

  • The element-wise worker directly called the private Arrow-to-pandas array converter, coupling it to column-level conversion internals.
  • A RecordBatch conversion boundary allows lower-level implementations to change without worker edits and avoids duplicating the PyArrow compatibility fallback.

Does this PR introduce any user-facing change?

No. This is an internal refactor that preserves existing UDF inputs, outputs, validation, and iterator behavior.

How was this patch tested?

  • Added test_concat_batches for multi-batch concatenation and the single-batch identity fast path.
  • Ran the conversion, higher-order-function UDF, and Arrow eval-handler suites under pandas 2/PyArrow 24 and pandas 3/PyArrow 25; both profiles passed 125 tests and 117 subtests. Ruff, formatting, and mypy also passed.

Was this patch authored or co-authored using generative AI tooling?

Generated-by: Codex (GPT-5)

@Spenserrrr
Spenserrrr force-pushed the arrow-worker-batch-interface branch from da35756 to adacdeb Compare September 18, 2026 23:23
@Spenserrrr
Spenserrrr marked this pull request as ready for review September 18, 2026 23:27
@Spenserrrr

Spenserrrr commented Sep 18, 2026 •

Copy link
Copy Markdown
Contributor Author

Hi @gaogaotiantian @Yicong-Huang! When I worked on #58751, I noticed that the element-wise read path still handled per-array conversion directly in worker.py, so I opened this PR. It moves the flatten/re-nest operations behind ArrowBatchTransformer and uses the batch-level ArrowToPandasConversion.to_pandas API. Could you take a look when you have time and let me know whether this boundary matches what you had in mind? Thanks!

@gaogaotiantian

Copy link
Copy Markdown
Contributor

I read the code and I think the direction is correct. We unified the logic to flat and reconstruct the nested data and ArrowBatchTransformer looks like the correct place to hold the new methods. @Yicong-Huang has a better understanding of what's going on for these specific transforms. I have no objection at this point.

@uros-b uros-b left a comment

Copy link
Copy Markdown
Member

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

+1 on the direction, leaving up to expert Python committers for further review, e.g. @Yicong-Huang

@Yicong-Huang Yicong-Huang left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Thanks @Spenserrrr! I support the idea to let eval types handle batch level logic only: this is a better abstraction so that eval types don't need to worry about low level implementation, but only focus on what kind of transformation on a batch is conducted. This PR aligns with my previous refactoring, but some details to be discussed. I left some comments inline.

Comment thread python/pyspark/sql/conversion.py Outdated
Comment thread python/pyspark/sql/conversion.py Outdated
Comment thread python/pyspark/sql/conversion.py Outdated
Comment thread python/pyspark/sql/conversion.py Outdated
Comment thread python/pyspark/sql/conversion.py Outdated
@Spenserrrr

Copy link
Copy Markdown
Contributor Author

Hi @gaogaotiantian @Yicong-Huang @uros-b, thanks for the review! I just pushed a revision addressing the review comments. I kept concat_batches general and moved the element-wise-specific shape helpers back to worker.py. Could you take another pass when you have time? Thanks!

Comment thread python/pyspark/worker.py Outdated

@Yicong-Huang Yicong-Huang left a comment

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

LGTM, thanks!

)

@classmethod
def concat_batches(cls, batches: Iterable["pa.RecordBatch"]) -> "pa.RecordBatch":

Copy link
Copy Markdown
Contributor

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

the only remaining concern is that if we need to do some validity checks

  • are all batches sharing the same schema
  • will the total size exceeding the upper limit of a record batch (I recall 2GB). otherwise it would need chunked arrays.

But since this is the current behavior we can keep it this way and rely on arrow's error messages for now.

Copy link
Copy Markdown
Contributor Author

Choose a reason for hiding this comment

The reason will be displayed to describe this comment to others. Learn more.

Agreed. I checked that both pa.concat_batches and the pre-19 pa.concat_arrays fallback already reject incompatible schemas. The size limit depends on the Arrow type and buffer layout, so a generic precheck here would be incomplete. I’ll keep the current behavior and let Arrow report these errors for now.

@Spenserrrr
Spenserrrr force-pushed the arrow-worker-batch-interface branch from 9e4ddae to b238c0e Compare September 23, 2026 04:46
@Spenserrrr

Copy link
Copy Markdown
Contributor Author

/trigger all

@Spenserrrr
Spenserrrr force-pushed the arrow-worker-batch-interface branch from b238c0e to 39f77b7 Compare September 23, 2026 18:32
Yicong-Huang pushed a commit that referenced this pull request Sep 23, 2026
…atch interfaces

### What changes were proposed in this pull request?

- Add a general `ArrowBatchTransformer.concat_batches` classmethod with compatibility for PyArrow versions before 19.0.0.
- Keep element-wise list flattening and shape reconstruction local to the worker while wrapping flattened inputs in RecordBatches for batch-level pandas/Arrow conversion.
- Reuse the shared concat helper in the grouped aggregate and window Arrow handlers migrated by [#58956](#58956), while preserving existing UDF contracts.

### Why are the changes needed?

- The element-wise worker directly called the private Arrow-to-pandas array converter, coupling it to column-level conversion internals.
- A RecordBatch conversion boundary allows lower-level implementations to change without worker edits and avoids duplicating the PyArrow compatibility fallback.

### Does this PR introduce _any_ user-facing change?

No. This is an internal refactor that preserves existing UDF inputs, outputs, validation, and iterator behavior.

### How was this patch tested?

- Added `test_concat_batches` for multi-batch concatenation and the single-batch identity fast path.
- Ran the conversion, higher-order-function UDF, and Arrow eval-handler suites under pandas 2/PyArrow 24 and pandas 3/PyArrow 25; both profiles passed 125 tests and 117 subtests. Ruff, formatting, and mypy also passed.

### Was this patch authored or co-authored using generative AI tooling?

Generated-by: Codex (GPT-5)

Closes #58903 from Spenserrrr/arrow-worker-batch-interface.

Authored-by: Spenser Sun <hsun112358@gmail.com>
Signed-off-by: Yicong-Huang <17627829+Yicong-Huang@users.noreply.github.com>
(cherry picked from commit 89062b8)
Signed-off-by: Yicong-Huang <17627829+Yicong-Huang@users.noreply.github.com>
@Yicong-Huang

Copy link
Copy Markdown
Contributor

Merge Summary:

Posted by merge_spark_pr.py

Sign up for free to join this conversation on GitHub. Already have an account? Sign in to comment

Labels

None yet

Projects

None yet

Development

Successfully merging this pull request may close these issues.

5 participants